DequeueType

Dequeue

Dequeue<'a> is the consuming side of a queue. Every Queue is a Dequeue, and Hub.subscribe returns a bare Dequeue, so one consumer works for both, and a subscriber cannot offer into its own buffer.

take waits for the next value; takeBetween waits for a minimum and takes a batch; poll, takeUpTo, and takeAll never wait. Values leave in FIFO order even when takers are interrupted: waiting takers are served one at a time, so a value given back by an interrupted taker returns before anything later is taken.

shutdown ends a queue without discarding its backlog: later takes drain what remains, and on a hub subscription it unsubscribes. stats reads the backlog and the accepted, dropped, and evicted counters for monitoring.

Consume a queue as a stream with FlowStream.fromDequeue. The Queue guide covers batching, shutdown, and monitoring.

Specification

Kind
Type
Members
11
Examples
0
takeSignature
Dequeue.take queue
pollSignature
Dequeue.poll queue
takeUpToSignature
Dequeue.takeUpTo max queue
takeBetweenSignature
Dequeue.takeBetween min max queue
takeAllSignature
Dequeue.takeAll queue

Summary

NameSignatureSynopsis
Operations
takeDequeue.take queueRemoves the oldest value, suspending until one is available.
pollDequeue.poll queueRemoves the oldest value if one is available, without suspending.
takeUpToDequeue.takeUpTo max queueRemoves up to max available values in FIFO order, without suspending.
takeBetweenDequeue.takeBetween min max queueRemoves between min and max values in FIFO order, suspending until at least min are available.
takeAllDequeue.takeAll queueRemoves every available value in FIFO order, without suspending.
sizeDequeue.size queueReturns the number of buffered values.
statsDequeue.stats queueReads the queue's backlog and counters in one consistent snapshot.
capacityDequeue.capacity queueReturns the queue's capacity, or None for an unbounded queue.
shutdownDequeue.shutdown queueShuts the queue down.
isShutdownDequeue.isShutdown queueReturns whether the queue has been shut down.
awaitShutdownDequeue.awaitShutdown queueSuspends until the queue is shut down.

Operations

kind:member

take

Dequeue.take queue
Member
Removes the oldest value, suspending until one is available.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, 'a>

Verification Examples

flow {
    let! (jobs: Queue<string>) = Queue.bounded 8
    do! jobs |> Queue.offer "job" |> Flow.ignore
    return! Dequeue.take jobs
}
kind:member

poll

Dequeue.poll queue
Member
Removes the oldest value if one is available, without suspending.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, 'a option>
kind:member

takeUpTo

Dequeue.takeUpTo max queue
Member
Removes up to max available values in FIFO order, without suspending.

Parameters

NameTypeDescription
maxint
queueDequeue<'a>

Returns

Flow<'env, 'error, 'a list>
kind:member

takeBetween

Dequeue.takeBetween min max queue
Member
Removes between min and max values in FIFO order, suspending until at least min are available.

Parameters

NameTypeDescription
minint
maxint
queueDequeue<'a>

Returns

Flow<'env, 'error, 'a list>

Verification Examples

// Wait for at least one sample, then write up to 500 in one batch.
samples |> Dequeue.takeBetween 1 500
kind:member

takeAll

Dequeue.takeAll queue
Member
Removes every available value in FIFO order, without suspending.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, 'a list>
kind:member

size

Dequeue.size queue
Member
Returns the number of buffered values.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, int>
kind:member

stats

Dequeue.stats queue
Member
Reads the queue's backlog and counters in one consistent snapshot.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, QueueStats>

Verification Examples

flow {
    let! stats = Dequeue.stats historian
    if stats.Capacity |> Option.exists (fun limit -> stats.Size * 10 >= limit * 8) then
        do! raiseAlarm "historian buffer 80% full"
}
kind:member

capacity

Dequeue.capacity queue
Member
Returns the queue's capacity, or None for an unbounded queue.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

int option
kind:member

shutdown

Dequeue.shutdown queue
Member
Shuts the queue down. Calling it again has no effect.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, unit>
kind:member

isShutdown

Dequeue.isShutdown queue
Member
Returns whether the queue has been shut down.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, bool>
kind:member

awaitShutdown

Dequeue.awaitShutdown queue
Member
Suspends until the queue is shut down.

Parameters

NameTypeDescription
queueDequeue<'a>

Returns

Flow<'env, 'error, unit>